Skip to content

HIVE-27126: queue level resource stats for YARN RM. - #6501

Open
architjainjain wants to merge 1 commit into
apache:masterfrom
architjainjain:HIVE-27126-yarn-RM
Open

HIVE-27126: queue level resource stats for YARN RM.#6501
architjainjain wants to merge 1 commit into
apache:masterfrom
architjainjain:HIVE-27126-yarn-RM

Conversation

@architjainjain

@architjainjain architjainjain commented May 20, 2026

Copy link
Copy Markdown
Contributor

What changes were proposed in this pull request?

This PR adds real-time YARN queue resource metrics display alongside Tez job progress during query execution. Users can now see memory, vCores, capacity utilization, running/pending applications, and allocated/pending containers for the queue being used by their queries.

Key Components:

  • YarnQueueMetricsCollector: Collects queue metrics from YARN ResourceManager
  • QueueMetricsCache: JVM-wide shared cache with expiry-based refresh (fetches only when cached data expires)
  • QueueMetricsRefreshPool: Scheduled executor pool for non-blocking metric collection
  • TezProgressMonitor: Enhanced to display queue metrics alongside task progress

What is NOT Changed:

  • Tez query execution logic remains unchanged
  • No modifications to YARN ResourceManager communication protocols
  • DAG submission, task scheduling, and container allocation logic untouched
  • Existing progress monitoring and logging infrastructure preserved

Technical Highlights:

  • Queue-level RM API calls: Single ResourceManager call per queue (not per session), only when cache expires
  • Shared caching: Multiple sessions using the same queue share cached metrics, reducing RM load
  • Dynamic refresh scheduling: Cache automatically adjusts to the lowest refresh interval among all active sessions for a queue
  • Jitter: Random 0-20% delay added to refresh intervals to prevent thundering herd when multiple queues expire simultaneously
  • Circuit breaker: Automatic fallback to no-op collector after 3 consecutive RM failures, preventing cascading failures

Configuration:

  • hive.tez.queue.metrics.refresh.interval: Per-session refresh interval (default: 10s, set to 0s to disable)
  • hive.server2.tez.queue.metrics.refresh.threads: Pool size for metric collection (default: 4)

Important Note: An HTTP exporter and Prometheus/Grafana integration were implemented as a local testing setup to validate the feature during development. These components are NOT part of this PR, NOT included in production code, and were only used in a local Docker environment for validation purposes.

Why are the changes needed?

Currently, when a Hive query is slow or stalled, users cannot determine if the issue is due to insufficient queue resources or other factors. They see Tez task progress (e.g., "Map 1: 3/10") but have no visibility into whether the queue has available memory/vCores to execute tasks in parallel.

This enhancement provides real-time queue resource information, enabling users to:

  • Identify resource bottlenecks immediately (e.g., queue at 95% memory usage)
  • Distinguish between slow queries due to lack of resources vs. data processing overhead
  • Make informed decisions about queue selection or query timing
  • Understand if pending containers are waiting due to queue capacity limits

Backward Compatibility:

  • Fully backward compatible - feature is opt-in via configuration (disabled by default in this implementation, can be enabled by setting refresh interval)
  • When disabled (interval=0s), behavior is identical to previous Hive versions
  • No changes to existing APIs, query execution, or Tez integration
  • No impact on existing queries, scripts, or workflows

Performance Impact:

  • Minimal overhead: Shared cache with single RM call per queue (not per session)
  • Non-blocking: Metrics collected in background thread pool
  • Graceful degradation: Circuit breaker prevents cascading failures
  • Zero impact when disabled

Does this PR introduce any user-facing change?

Yes. When hive.tez.queue.metrics.refresh.interval is set to a positive value (default: 10s), users will see queue-level metrics displayed with Tez job progress:

In-place mode (hive.server2.in.place.progress=true):

Map 1: 3/10   Reducer 2: 0/5
QUEUE: default
MEMORY: 1.5/5.4 GB (27.78% used) | VCORES: 3/7 (42.86% used)
CAPACITY: 46.30% (used), 60.00% (allocated)
APPS: 1 running, 0 pending | CONTAINERS: 4 allocated, 0 pending

Log file mode (hive.server2.in.place.progress=false):

INFO  : Map 1: 3/10   Reducer 2: 0/5    QUEUE: default | MEMORY: 1.5/5.4 GB (27.78% used) | VCORES: 3/7 (42.86% used) | CAPACITY: 46.30% (used), 60.00% (allocated) | APPS: 1 running, 0 pending | CONTAINERS: 4 allocated, 0 pending

When disabled (set hive.tez.queue.metrics.refresh.interval=0s):
No queue metrics are displayed, behavior remains identical to previous versions.

How was this patch tested?

The patch was tested in a multi-node YARN cluster environment with concurrent queries running on different queues.

Testing included:

  1. Functional testing: Queue metrics displayed correctly with various refresh intervals (1s, 5s, 10s)
  2. Feature toggle testing: Metrics disabled when refresh interval set to 0s
  3. Multi-session testing: Multiple concurrent sessions with different queues and refresh intervals
  4. Dynamic scheduling: Verified cache adjusts to lowest interval when sessions with different settings share a queue
  5. Cross-mode testing: Verified consistent output in both in-place and log file modes
  6. Performance testing: Validated minimal overhead with shared cache and background refresh pool
  7. Efficiency validation: Confirmed single RM API call per queue (not per session) only on cache expiry
  8. Jitter behavior: Verified 0-20% randomization prevents simultaneous refresh across queues
  9. Circuit breaker: Validated automatic failover to no-op collector after consecutive RM failures

Test Environment:

  • YARN cluster with multiple queues (default, analytics, batch)
  • HiveServer2 with Tez execution engine
  • Concurrent query workloads

Validation artifacts (local testing setup only - NOT in production):

  • Prometheus exporter and Grafana dashboards were implemented in a local Docker environment to validate metric accuracy, cache behavior, and refresh scheduling
  • Multi-session test scripts for concurrent scenario validation
  • These testing tools are separate from the production code and not included in this PR

Screenshots:

1. Terminal output with queue metrics enabled (hive.tez.queue.metrics.refresh.interval=10s):

Queue metrics showing memory, vcores, capacity, apps and containers alongside Tez progress

2. Terminal output with queue metrics disabled (hive.tez.queue.metrics.refresh.interval=0s):

Query progress without queue metrics, showing only Tez task progress

3. Grafana dashboard showing backend metrics validation (local testing setup only):

Screenshot 2026-07-22 at 11 29 39 AM

Note: The Grafana and Prometheus infrastructure shown above was created exclusively for local testing and validation purposes. This monitoring stack is NOT part of the production code, NOT included in the deployment, and was only used in a local Docker environment to validate the feature's correctness.

Configuration tested:

<!-- Enable queue metrics with 10-second refresh -->
<property>
  <name>hive.tez.queue.metrics.refresh.interval</name>
  <value>10s</value>
</property>

<!-- Disable queue metrics -->
<property>
  <name>hive.tez.queue.metrics.refresh.interval</name>
  <value>0s</value>
</property>

@abstractdog abstractdog left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

thanks @architjainjain so far, dropped some comments, but I wasn't able to fully read it through, I'll get back after, in the meantime I can do some testing too hopefully

* behaviour added as part of HIVE-27126.
*
* We capture stdout via a ByteArrayOutputStream and inspect the rendered output.
*/public class TestInPlaceUpdate {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

line break before public

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

done

Comment thread common/src/test/org/apache/hadoop/hive/common/log/TestInPlaceUpdate.java Outdated
Comment thread common/src/test/org/apache/hadoop/hive/common/log/TestInPlaceUpdate.java Outdated
Comment thread common/src/test/org/apache/hadoop/hive/common/log/TestInPlaceUpdate.java Outdated
Comment thread common/src/test/org/apache/hadoop/hive/common/log/TestInPlaceUpdate.java Outdated
Comment on lines +318 to +324
for (int i = 0; i < 5; i++) {
collectors[i] = new YarnQueueMetricsCollector(
mockYarnClient, "default", refreshIntervalMs, "jitter-test-query-" + i);
assertNotNull("Collector " + i + " should be created successfully", collectors[i]);
}
// If we get here, all 5 collectors were created with their own jittered delays
// without conflict or exception - thundering herd fix is in place

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

considering mock yarn clients, creating concurrent collectors for jitter testing doesn't seem that valuable to me...I don't think they will ever throw an exception

mockYarnClient, "default", 30, "circuit-breaker-recovery-query");

try {
// Wait for recovery - snapshot should eventually be populated

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

need extra assertion about the failure first to make 100% sure we actually hit circuit breaker logic

Comment on lines +428 to +429
public void testCircuitBreakerDoesNotAffectSuccessfulCollection() throws Exception {
// Normal operation - no failures, circuit breaker should never activate

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this is just a simple happy testing, I guess it's covered by testSuccessfulMetricsCollection

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@abstractdog Thanks for catching this! I've added JavaDoc to clarify that this test is part of the circuit breaker test suite. While it looks similar to testSuccessfulMetricsCollection(), it specifically validates that the circuit breaker doesn't interfere with normal operations. The circuit breaker suite needs three scenarios: failure, recovery, and happy path. This test covers the happy path behavior for circuit breaker specifically, not general metrics collection. I've added cross-references to the related tests for clarity.

Comment on lines +485 to +493
when(mockQueueStats.getAllocatedMemoryMB()).thenReturn(1024L);
when(mockQueueStats.getAvailableMemoryMB()).thenReturn(1024L);
when(mockQueueStats.getAllocatedVCores()).thenReturn(4L);
when(mockQueueStats.getAvailableVCores()).thenReturn(4L);
when(mockQueueStats.getNumAppsRunning()).thenReturn(1L);
when(mockQueueStats.getPendingContainers()).thenReturn(0L);
when(mockQueueInfo.getQueueStatistics()).thenReturn(mockQueueStats);
when(mockQueueInfo.getCapacity()).thenReturn(0.5f);
when(mockYarnClient.getQueueInfo("default")).thenReturn(mockQueueInfo);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I feel like defining all these variables adds a lot of boilerplate to the unit tests, whereas in this case we're only interested in collection timestamp, could a separate method be created just for adding happy value, and it can be overridden per test, like:

  when(mockQueueStats.getAllocatedMemoryMB()).thenReturn(1024L);
  when(mockQueueStats.getAvailableMemoryMB()).thenReturn(1024L);
  when(mockQueueStats.getAllocatedVCores()).thenReturn(4L);
  when(mockQueueStats.getAvailableVCores()).thenReturn(4L);
  when(mockQueueStats.getNumAppsRunning()).thenReturn(1L);
  when(mockQueueStats.getPendingContainers()).thenReturn(0L);
  when(mockQueueInfo.getQueueStatistics()).thenReturn(mockQueueStats);
  when(mockQueueInfo.getCapacity()).thenReturn(0.5f);
  when(mockYarnClient.getQueueInfo("default")).thenReturn(mockQueueInfo);

method javadoc can also tell that these are only default happy values for the unit tests

Comment on lines +824 to +830
// HIVE-27126: Thrift regeneration omitted setStartTimeIsSet(true) from the constructor.
// Explicitly call setStartTime() to set the isset flag required for Thrift validation.
tProgressUpdateResp.setStartTime(progressUpdate.startTimeMillis);
if (progressUpdate.queueMetrics() != null && !progressUpdate.queueMetrics().isEmpty()) {
tProgressUpdateResp.setQueueMetrics(progressUpdate.queueMetrics());
}
resp.setProgressUpdateResponse(tProgressUpdateResp);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

why is this change needed?

@architjainjain architjainjain Jun 23, 2026

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This works around a Thrift code generation bug. When we regenerated Thrift code after adding the queueMetrics field, the generated constructor for TProgressUpdateResp accepts startTimeMillis but fails to set the isset flag. Without this flag, Thrift serialization treats the field as unset, causing clients to receive incomplete progress updates. The explicit setStartTime() call ensures both the value AND the isset flag are properly set, maintaining backward compatibility with Thrift clients.

without this test the generated file is not having the starttime isset flag updated.

https://github.com/apache/hive/pull/6501/changes#diff-d71ddeafeea57e0fbf6a9dfda436eab06b69a39ac46c815cc1b15bfffe5f35e0L155

this.status = status;
this.footerSummary = footerSummary;
this.startTime = startTime;
setStartTimeIsSet(true);

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@abstractdog this is getting removed. so we added manually after the constructor call.

@sonarqubecloud

Copy link
Copy Markdown

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Note

Copilot couldn't run its full agentic review because it didn't start before the timeout. Make sure your repository has a runner available, or add a copilot-code-review.yml file specifying one with the runs-on attribute. See the docs for more details.

Adds optional, real-time YARN queue resource metrics to the Tez progress display during Hive query execution, propagating the data through HS2/Thrift and Beeline in-place updates.

Changes:

  • Introduces YARN queue metrics collection (collector + per-queue state + shared cache + scheduled refresh pool with jitter/circuit breaker).
  • Integrates queue metrics into Tez progress rendering (in-place + log-to-file) and publishes it over Thrift (TProgressUpdateResp).
  • Adds unit tests covering cache/state/pool behavior and Tez monitor formatting/initialization.

Reviewed changes

Copilot reviewed 29 out of 30 changed files in this pull request and generated 9 comments.

Show a summary per file
File Description
service/src/java/org/apache/hive/service/server/HiveServer2.java Initializes the shared refresh pool during HS2 Tez session pool startup.
service/src/java/org/apache/hive/service/cli/thrift/ThriftCLIService.java Populates new Thrift queueMetrics field and applies Thrift “isset” workaround for startTime.
service/src/java/org/apache/hive/service/cli/JobProgressUpdate.java Adds queue metrics string to progress update model.
service-rpc/if/TCLIService.thrift Adds optional queueMetrics to TProgressUpdateResp.
ql/src/java/org/apache/hadoop/hive/ql/session/SessionState.java Implements queueMetrics() default for the progress monitor facade.
ql/src/java/org/apache/hadoop/hive/ql/exec/tez/monitoring/yarnqueue/QueueMetricsCollector.java New interface for queue metrics collectors.
ql/src/java/org/apache/hadoop/hive/ql/exec/tez/monitoring/yarnqueue/NoOpQueueMetricsCollector.java Null-object collector when feature is disabled.
ql/src/java/org/apache/hadoop/hive/ql/exec/tez/monitoring/yarnqueue/QueueMetricsSnapshot.java Immutable snapshot of queue resource stats fetched from RM.
ql/src/java/org/apache/hadoop/hive/ql/exec/tez/monitoring/yarnqueue/QueueMetricsState.java Per-queue state: interval registration, scheduling, refresh lock, circuit breaker.
ql/src/java/org/apache/hadoop/hive/ql/exec/tez/monitoring/yarnqueue/QueueMetricsCache.java JVM-wide Guava cache mapping queue name → state with expiry.
ql/src/java/org/apache/hadoop/hive/ql/exec/tez/monitoring/yarnqueue/QueueMetricsRefreshPool.java JVM-wide scheduled executor singleton to run refresh tasks (with jitter).
ql/src/java/org/apache/hadoop/hive/ql/exec/tez/monitoring/yarnqueue/YarnQueueMetricsCollector.java Active collector that coordinates with shared cache/state and refresh pool.
ql/src/java/org/apache/hadoop/hive/ql/exec/tez/monitoring/TezProgressMonitor.java Formats and returns multi-line queue metrics for progress rendering.
ql/src/java/org/apache/hadoop/hive/ql/exec/tez/monitoring/TezJobMonitor.java Creates/shuts down the metrics collector based on config and wires it into progress monitor.
ql/src/java/org/apache/hadoop/hive/ql/exec/tez/monitoring/RenderStrategy.java Appends queue metrics into log-to-file progress report (single-line rendering).
ql/src/java/org/apache/hadoop/hive/ql/exec/tez/TezSessionState.java Adds per-session YarnClient lifecycle to support collector RM queries.
ql/src/java/org/apache/hadoop/hive/ql/exec/tez/TezSessionPoolSession.java Exposes YarnClient via pooled session wrapper.
ql/src/java/org/apache/hadoop/hive/ql/exec/tez/TezSession.java Extends TezSession interface with getYarnClient().
common/src/java/org/apache/hadoop/hive/conf/HiveConf.java Adds new configuration keys for refresh interval and refresh thread pool size.
common/src/java/org/apache/hadoop/hive/common/log/ProgressMonitor.java Adds queueMetrics() to the ProgressMonitor contract.
common/src/java/org/apache/hadoop/hive/common/log/InPlaceUpdate.java Renders queue metrics block (multi-line) under the standard progress output.
beeline/src/java/org/apache/hive/beeline/logs/BeelineInPlaceUpdateStream.java Reads queueMetrics from Thrift progress updates for in-place rendering.
ql/src/test/org/apache/hadoop/hive/ql/exec/tez/monitoring/yarnqueue/TestQueueMetricsState.java Unit tests for per-queue state logic (intervals, refresh lock, circuit breaker).
ql/src/test/org/apache/hadoop/hive/ql/exec/tez/monitoring/TestTezProgressMonitorQueueMetrics.java Unit tests for queue metrics formatting and edge cases in TezProgressMonitor.
ql/src/test/org/apache/hadoop/hive/ql/exec/tez/monitoring/TestTezJobMonitorQueueMetrics.java Unit tests for TezJobMonitor collector initialization decisions.
ql/src/test/org/apache/hadoop/hive/ql/exec/tez/TestYarnQueueMetricsCollector.java Unit tests for collector/cache/pool integration behavior and circuit breaker.
ql/src/test/org/apache/hadoop/hive/ql/exec/tez/TestQueueMetricsRefreshPool.java Unit tests for refresh pool singleton, scheduling, and jitter determinism/range.
ql/src/test/org/apache/hadoop/hive/ql/exec/tez/TestQueueMetricsCache.java Unit tests for cache placeholder/put semantics and concurrency behavior.
ql/src/test/org/apache/hadoop/hive/ql/exec/tez/TestNoOpQueueMetricsCollector.java Unit tests for NoOp collector behavior and singleton semantics.
Files not reviewed (1)
  • service-rpc/src/gen/thrift/gen-javabean/org/apache/hive/service/rpc/thrift/TProgressUpdateResp.java: Generated file

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment on lines +815 to +825
if (session != null && yarnClient == null) {
try {
yarnClient = YarnClient.createYarnClient();
yarnClient.init(conf);
yarnClient.start();
LOG.info("YarnClient initialized for session: {}", sessionId);
} catch (Exception e) {
LOG.warn("Failed to initialize YarnClient for metrics collection", e);
yarnClient = null;
}
}
Comment on lines +125 to +128
public static long calculateJitter(String queueName, long intervalMs) {
long jitterWindow = intervalMs * JITTER_PERCENT / 100;
return Math.abs(queueName.hashCode()) % jitterWindow;
}
Comment on lines +113 to +115
public ScheduledFuture<Void> scheduleRefreshTask(Runnable task, long intervalMs) {
return (ScheduledFuture<Void>) refreshPool.scheduleWithFixedDelay(task, 0, intervalMs, TimeUnit.MILLISECONDS);
}
Comment on lines +973 to +977
/**
* Initializes the queue metrics refresh pool and HTTP exporter.
* Failures are non-fatal — logged as warnings so the server can start without queue metrics.
*/
private void initializeQueueMetricsPool(HiveConf hiveConf) {
Comment on lines +982 to +984
} catch (Exception e) {
LOG.warn("Failed to initialize queue metrics refresh pool: {}", e.getMessage());
}
private static final Logger LOG = LoggerFactory.getLogger(QueueMetricsRefreshPool.class);

private static final int DEFAULT_THREAD_COUNT = 4;
public static final int JITTER_PERCENT = 10;
Comment on lines +4009 to +4014
HIVE_TEZ_QUEUE_METRICS_REFRESH_INTERVAL("hive.tez.queue.metrics.refresh.interval", "0s",
new TimeValidator(TimeUnit.SECONDS),
"Interval for refreshing YARN queue resource metrics during Tez query execution. " +
"When set to a positive value (e.g. 10s), displays real-time memory, vCore, capacity " +
"and application metrics for the YARN queue being used. " +
"Set to 0 or negative to disable. Minimum effective value is 1 second."),
Comment on lines +95 to +102
long startTime = System.currentTimeMillis();
QueueMetricsSnapshot snapshot;
while ((snapshot = collector.getLatestSnapshot()) == null) {
if (System.currentTimeMillis() - startTime > timeoutMs) {
fail("Snapshot not available after " + timeoutMs + "ms");
}
}
return snapshot;
Comment on lines +220 to +227
for (Object[] row : cases) {
String label = (String) row[0];
YarnClient yarnClient = (YarnClient) row[1];
String queueName = (String) row[2];

// Re-initialise mocks so state from the previous row does not leak.
MockitoAnnotations.openMocks(this);
when(mockDag.getName()).thenReturn("test-dag-1");
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants